welo git mirrors
Commit 7665d1ff7cd726ad6916167d1bb693ac95d4e5e2
Parents : 97bc2fc
Author : welo | main nixos computer <empty@empty.com>
Date : 2026-07-09T06:25:13+01:00
cleaned up warnings and removed more print and info statements
Changes
8 files changed, 63 insertions(+), 84 deletions(-)
Diff
diff --git a/examples/server.rs b/examples/server.rs
index d1aa40c..95779ea 100644
--- a/examples/server.rs
+++ b/examples/server.rs
@@ -32,6 +32,7 @@ async fn main() {
match cli.command {
Commands::Server { identity_file } => {
+ drop(identity_file);
eprintln!("Starting RNS SOCKS5 proxy server...");
eprintln!("Make sure rnsd is running (pip install rns && rnsd)");
eprintln!();
diff --git a/src/client.rs b/src/client.rs
index 1b87cd5..51a9a05 100644
--- a/src/client.rs
+++ b/src/client.rs
@@ -16,21 +16,17 @@ use std::net::{IpAddr, Ipv4Addr, SocketAddr};
use std::sync::Arc;
use std::time::Duration;
-use clap::Command;
use fast_socks5::server::Socks5ServerProtocol;
-use fast_socks5::server::states::{CommandRead, Opened};
+use fast_socks5::server::states::CommandRead;
use fast_socks5::util::target_addr::TargetAddr;
use fast_socks5::{ReplyError, Socks5Command};
use log::{debug, error, info, warn};
-use rns_net::discovery::filter_and_sort_interfaces;
use rns_net::{LinkId, RnsNode};
use tokio::net::{TcpListener, TcpStream, UdpSocket};
use tokio::sync::mpsc::UnboundedReceiver;
use tokio::sync::{mpsc, Notify};
-use tokio::time::error::Elapsed;
-use crate::filter::filter_and_convert;
-use crate::forwarding::PortType::{self, Tcp};
+use crate::forwarding::PortType;
use crate::forwarding::{ForwardedPort, tcp_tunnel, udp_tunnel};
use crate::mux::MuxHandle;
use crate::relay::relay_bidirectional_udp_client_side;
@@ -50,14 +46,16 @@ pub async fn run_client(server_hex: &str, listen_addr: &str) {
pub async fn run_client_forward(server_hex: &str, ports: Vec<ForwardedPort> ) {
if let Some(destination) = decode_hash(server_hex).await {
if let Some((mux, node, rx)) = connect_rns(destination).await {
- info!("a");
let reconnect_notify = reconnect_generator(mux.clone(), node, rx, destination).await;
- info!("b");
for port in ports {
match port.r#type {
PortType::Tcp => {
- tcp_tunnel(mux.clone(), reconnect_notify.clone(), port).await
+ let notify = reconnect_notify.clone();
+ let mux = mux.clone();
+ tokio::spawn(async move {
+ tcp_tunnel(mux, notify, port).await
+ });
},
PortType::Udp => {
info!("c");
@@ -69,9 +67,10 @@ pub async fn run_client_forward(server_hex: &str, ports: Vec<ForwardedPort> ) {
}
}
}
+ let __ : () =std::future::pending().await; // so main thread never ends.
// run_sockets_proxy_handling(listen_addr, mux.clone(), reconnect_notify).await;
}
- }
+ };
}
pub async fn decode_hash(server_hex: &str) -> Option<[u8; 16]> {
@@ -263,9 +262,8 @@ async fn dispatch_and_reconnect(
match event {
ProxyEvent::LinkData { data, .. } => {
- // println!("{:?}",data);
for frame in mux.receive_data(&data).await {
- // println!("type {:?} sid: {:?}", frame.frame_type, frame.session_id);
+
mux.dispatch(frame).await;
}
}
@@ -363,7 +361,7 @@ pub async fn connect_tcp_server_side(
mux: MuxHandle,
session_rx: &mut mpsc::UnboundedReceiver<Frame>,
target_addr: TargetAddr,)
- -> Option<Result<(),String>>
+ -> Result<(),String>
{
//
// if let Some(socket_addr) = filter_and_convert(target_addr, None).await {
@@ -377,8 +375,6 @@ pub async fn connect_tcp_server_side(
// Wait for CONN_OK or CONN_ERR with timeout
tokio::time::timeout(Duration::from_secs(15), async {
while let Some(frame) = session_rx.recv().await {
- // println!("type: {:?}", frame.frame_type);
- // println!("type: {:?}", frame.session_id);
match frame.frame_type {
FrameType::ConnectOk => return Ok(()),
FrameType::ConnectErr => {
@@ -389,8 +385,10 @@ pub async fn connect_tcp_server_side(
}
}
Err("channel closed".to_string())
- })
- .await.ok()
+ }).await.map_err(|_e| "Timeout".into()).flatten()
+
+
+ // a.ok()
// } else {
// return None;
@@ -414,13 +412,13 @@ async fn handle_tcp_connect(
info!("[{}] -> {}:{} connection result {:?}", sid, host, port, connect_result);
- if let Some(_) = connect_result {
+ if let Ok(_) = connect_result {
// Reply to SOCKS5 client based on RNS connection result
let dummy_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(0, 0, 0, 0)), 0);
let stream = match connect_result {
- Some(Ok(())) => {
+ Ok(()) => {
// Connection succeeded -- send SOCKS5 success reply
match proto.reply_success(dummy_addr).await {
Ok(s) => s,
@@ -432,18 +430,12 @@ async fn handle_tcp_connect(
}
}
}
- Some(Err(reason)) => {
+ Err(reason) => {
warn!("[{}] Remote connect failed: {}", sid, reason);
let _ = proto.reply_error(&ReplyError::GeneralFailure).await;
mux.drop_session(sid).await;
return None;
}
- None => {
- warn!("[{}] Connect timeout", sid);
- let _ = proto.reply_error(&ReplyError::TtlExpired).await;
- mux.drop_session(sid).await;
- return None;
- }
};
@@ -464,7 +456,7 @@ pub async fn udp_bind_connect(
mux: MuxHandle,
session_rx: &mut mpsc::UnboundedReceiver<Frame>,
target_addr: TargetAddr,)
- -> Result<Result<(),String>,Elapsed>
+ -> Result<(),String>
{
let (host, port) = target_addr.into_string_and_port();
@@ -487,8 +479,7 @@ pub async fn udp_bind_connect(
}
}
Err("channel closed".to_string())
- })
- .await
+ }).await.map_err(|_e| "Timeout".into()).flatten()
}
async fn handle_udp_connect(
@@ -512,7 +503,7 @@ async fn handle_udp_connect(
match connect_result {
- Ok(Ok(())) => {
+ Ok(()) => {
// Connection succeeded -- send SOCKS5 success reply
match proto.reply_success(relay_address).await {
Ok(s) => Some((udp_stream,s)),
@@ -524,18 +515,13 @@ async fn handle_udp_connect(
}
}
}
- Ok(Err(reason)) => {
+ Err(reason) => {
warn!("[{}] Remote connect failed: {}", sid, reason);
+
let _ = proto.reply_error(&ReplyError::GeneralFailure).await;
mux.drop_session(sid).await;
return None;
}
- Err(_) => {
- warn!("[{}] Connect timeout", sid);
- let _ = proto.reply_error(&ReplyError::TtlExpired).await;
- mux.drop_session(sid).await;
- return None;
- }
}
diff --git a/src/forwarding.rs b/src/forwarding.rs
index ec5c55c..dc77032 100644
--- a/src/forwarding.rs
+++ b/src/forwarding.rs
@@ -22,7 +22,7 @@ use std::{net::{Ipv4Addr, SocketAddr}, sync::Arc};
use fast_socks5::util::target_addr::TargetAddr;
use log::{error, info, warn};
-use tokio::{net::{TcpListener, UdpSocket}, sync::{Notify, mpsc::{UnboundedReceiver, UnboundedSender, channel, unbounded_channel}}};
+use tokio::{net::{TcpListener}, sync::Notify, };
use udp_stream::UdpListener;
use crate::{client::{connect_tcp_server_side, udp_bind_connect}, mux::MuxHandle, relay::{relay_forwarded_tcp, relay_forwarded_udp}};
@@ -54,7 +54,7 @@ pub async fn tcp_tunnel(mux: MuxHandle, reconnect_notify: Arc<Notify> , port: Fo
let listener = match TcpListener::bind(format!("127.0.0.1:{}",port.client_port)).await {
Ok(l) => l,
Err(e) => {
- error!("Failed to local port at {}", port.client_port);
+ error!("Failed to local port at {} cause of {:?}", port.client_port, e);
return;
}
};
@@ -65,9 +65,9 @@ pub async fn tcp_tunnel(mux: MuxHandle, reconnect_notify: Arc<Notify> , port: Fo
loop {
tokio::select! {
accept_result = listener.accept() => {
- let (stream, addr) = match accept_result {
+ let (stream, _addr) = match accept_result {
Ok(sa) => sa,
- Err(e) => {
+ Err(_e) => {
continue;
}
};
@@ -84,8 +84,13 @@ pub async fn tcp_tunnel(mux: MuxHandle, reconnect_notify: Arc<Notify> , port: Fo
let target_addr = TargetAddr::Ip(SocketAddr::new(std::net::IpAddr::V4(Ipv4Addr::LOCALHOST), port.server_port)); // ask for any port client server side
tokio::spawn(async move {
- if let Some(_) = connect_tcp_server_side(sid, mux.clone(), &mut session_rx, target_addr ).await {
- relay_forwarded_tcp(sid, stream, mux_clone, session_rx).await;
+ match connect_tcp_server_side(sid, mux.clone(), &mut session_rx, target_addr ).await {
+ Ok(()) => {
+ relay_forwarded_tcp(sid, stream, mux_clone, session_rx).await;
+ }
+ Err(error_message) => {
+ warn!("[{}] could not connect forward cause of {}", sid, error_message);
+ }
}
});
}
@@ -108,48 +113,37 @@ pub async fn udp_tunnel(mux: MuxHandle, reconnect_notify: Arc<Notify> , port: Fo
let target_addr = SocketAddr::new(std::net::IpAddr::V4(Ipv4Addr::LOCALHOST), port.client_port);
let listener = match UdpListener::bind(target_addr).await {
Ok(l) => l,
- Err(e) => {
+ Err(_e) => {
error!("Failed to local port at {}", port.client_port);
return;
}
};
- info!("d");
-
let reference = &mux;
loop {
tokio::select! {
accept_result = listener.accept() => {
- let (stream, addr) = match accept_result {
+ let (stream, _addr) = match accept_result {
Ok(sa) => sa,
- Err(e) => {
- info!("hi");
+ Err(_e) => {
continue;
}
};
- info!("5");
let mux = reference.clone();
- info!("1");
if !mux.is_connected().await {
drop(stream);
continue;
}
- info!("2");
let sid = mux.next_session_id().await;
- info!("2");
let mut session_rx = mux.register_session(sid).await;
- info!("2");
let mux_clone = mux.clone();
- info!("2");
- info!("2");
// let connect_result = udp_bind_connect(sid,mux.clone(), &mut session_rx, target_addr).await;
let target_addr = TargetAddr::Ip(SocketAddr::new(std::net::IpAddr::V4(Ipv4Addr::LOCALHOST), port.server_port ));
- info!("3");
tokio::spawn(async move {
if let Ok(_) = udp_bind_connect(sid, mux.clone(), &mut session_rx, target_addr ).await {
relay_forwarded_udp(sid, stream, mux_clone, session_rx, port.server_port).await;
diff --git a/src/frame.rs b/src/frame.rs
index 512f7c7..cb7dd7a 100644
--- a/src/frame.rs
+++ b/src/frame.rs
@@ -9,8 +9,6 @@
use std::fmt;
-use clap::error::ErrorFormatter;
-use log::{error, info, warn};
use crate::frame::FrameDecodeState::{DecodingFailed, MoreDataRequired};
diff --git a/src/main.rs b/src/main.rs
index 6ab2204..3b9959d 100644
--- a/src/main.rs
+++ b/src/main.rs
@@ -47,7 +47,6 @@ async fn main() {
Commands::Connect { destination, tcp_port, udp_port, both_port} => {
let ports_to_connect = connect_vec_merging_connect(tcp_port, udp_port, both_port);
- println!("{:?}",ports_to_connect);
run_client_forward(&destination, ports_to_connect).await;
// rns_proxy::server::run_server(identity_file.as_deref()).await;
}
diff --git a/src/mux.rs b/src/mux.rs
index 567768e..4ce3fd8 100644
--- a/src/mux.rs
+++ b/src/mux.rs
@@ -15,14 +15,12 @@
//! in rns-rs, causing `NotReady` errors after the first 2 messages.
use std::collections::HashMap;
-use std::error;
use std::sync::Arc;
use tokio::sync::{Mutex};
-use log::{debug, error, info, warn};
+use log::{debug, error, warn};
use rns_core::constants::LINK_MDU;
-use rns_core::msgpack::Error;
-use rns_net::{LinkId, LocalServerFactory, RnsNode};
+use rns_net::{LinkId, RnsNode};
use tokio::sync::mpsc;
use crate::frame::FrameDecodeState::{DecodingFailed, MoreDataRequired};
@@ -41,7 +39,7 @@ pub struct MuxHandle {
struct MuxInner {
node: Arc<RnsNode>,
link_id: Mutex<Option<LinkId>>,
- sessions: Mutex<HashMap<u32, mpsc::UnboundedSender<Frame>>>,
+ sessions: Mutex<HashMap<u32, tokio::sync::mpsc::UnboundedSender<Frame>>>,
next_sid: Mutex<u32>,
/// Reassembly buffer for incoming raw link data chunks.
recv_buf: Mutex<Vec<u8>>,
@@ -87,7 +85,7 @@ impl MuxHandle {
/// Get the next session id.
pub async fn next_session_id(&self) -> u32 {
- let mut sid = (self.inner.next_sid.lock().await);
+ let mut sid = self.inner.next_sid.lock().await;
*sid = sid.wrapping_add(1);
*sid
}
@@ -152,8 +150,11 @@ impl MuxHandle {
let sessions = self.inner.sessions.lock().await;
if let Some(tx) = sessions.get(&sid) {
+ // info!("hi");
if tx.send(frame).is_err() {
- debug!("Session {} channel closed", sid);
+ // debug!("Session {} channel closed", sid);
+ } else {
+ // info!("fine?");
}
} else {
warn!("No session {} for frame type {}", sid, ft);
diff --git a/src/relay.rs b/src/relay.rs
index 634e9ab..bc1ff24 100644
--- a/src/relay.rs
+++ b/src/relay.rs
@@ -3,15 +3,13 @@
//! Note 1, idk where to put this but the buffer has to be bigger than 64k as v
//! orginally was 4096 but occasionally there would be packets
-use std::net::{Ipv4Addr, ToSocketAddrs};
-use std::os::unix::net::SocketAddr;
-use std::str::SplitWhitespace;
+use std::net::Ipv4Addr;
use std::sync::Arc;
-use std::time::Duration;
+
use fast_socks5::{new_udp_header, parse_udp_request};
-use fast_socks5::util::target_addr::{TargetAddr, ToTargetAddr};
-use log::{debug, error, info, warn};
+use fast_socks5::util::target_addr::ToTargetAddr;
+use log::{debug, warn};
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::sync::{Mutex, mpsc};
use udp_stream::UdpStream;
@@ -33,7 +31,7 @@ pub async fn relay_bidirectional_tcp(
let mux_fwd = mux.clone();
- info!("[{}] relay", sid);
+ debug!("[{}] relay bidirectional tcp", sid);
// TCP -> RNS
let tcp_to_rns = tokio::spawn(async move {
@@ -65,6 +63,7 @@ pub async fn relay_bidirectional_tcp(
}
}
FrameType::Close => break,
+ FrameType::ConnectErr => {warn!("connection errored, should probably be handled by different code"); break},
_ => {}
}
}
@@ -75,7 +74,7 @@ pub async fn relay_bidirectional_tcp(
_ = rns_to_tcp => {},
}
- info!("[{}] drop connection", sid);
+ debug!("[{}] drop connection", sid);
mux.send(FrameType::Close, sid, Vec::new()).await;
mux.drop_session(sid).await;
@@ -258,7 +257,7 @@ pub async fn relay_bidirectional_udp_server_side(
match frame.frame_type {
FrameType::Data => {
match parse_udp_request(&*frame.payload).await {
- Ok((frag,addr,data)) => {
+ Ok((_frag,addr,data)) => {
if let Some(socket) = filter_and_convert(addr.clone(), Some(&config2)).await {
@@ -327,6 +326,8 @@ pub async fn relay_forwarded_tcp(
let (mut tcp_read, mut tcp_write) = stream.into_split();
let mux_fwd = mux.clone();
+ debug!("relay forwarded tcp");
+
// TCP -> RNS
let tcp_to_rns = tokio::spawn(async move {
let mut buf = [0u8; 4096];
@@ -358,6 +359,7 @@ pub async fn relay_forwarded_tcp(
}
}
FrameType::Close => break,
+ FrameType::ConnectErr => {warn!("connection errored, should probably be handled by different code"); break},
_ => {}
}
}
@@ -368,6 +370,7 @@ pub async fn relay_forwarded_tcp(
_ = rns_to_tcp => {},
}
+ debug!("[{}] tcp relay ended", sid);
mux.send(FrameType::Close, sid, Vec::new()).await;
mux.drop_session(sid).await;
}
@@ -415,8 +418,8 @@ pub async fn relay_forwarded_udp(
FrameType::Data => {
match parse_udp_request(&*frame.payload).await {
- Ok((frag,addr,data)) => {
- let target = addr.into_string_and_port();
+ Ok((_frag,_addr,data)) => {
+ // let target = addr.into_string_and_port();
if let Err(e) = udp_write.write_all(data).await {
warn!("[{}] UDP write error: {}", sid, e);
@@ -447,6 +450,7 @@ pub async fn relay_forwarded_udp(
_ = rns_to_tcp => {},
}
+ debug!("[{}] udp relay ended", sid);
mux.send(FrameType::Close, sid, Vec::new()).await;
mux.drop_session(sid).await;
}
diff --git a/src/server.rs b/src/server.rs
index 68529e6..40c6c23 100644
--- a/src/server.rs
+++ b/src/server.rs
@@ -168,7 +168,6 @@ pub async fn run_server(identity_path: Option<&str>, filter_config: FilterConfig
FrameType::Connect => {
let sid = frame.session_id;
if let Some((host, port,udp)) = decode_connect_payload(&frame.payload) {
- info!("[{}] -> {}:{} {}", sid, host, port, udp);
let addr = TargetAddr::Domain(host, port);
let session_rx = mux.register_session(sid).await;
let mux_clone = mux.clone();
@@ -216,9 +215,7 @@ async fn handle_server_session_tcp(
filter_config: FilterConfig
) {
- info!("{} start", sid);
if let Some(socket) = filter_and_convert(addr.clone(), Some(&filter_config)).await {
- info!("{} hi", sid);
let stream = match TcpStream::connect(socket).await {
Ok(s) => s,
Err(e) => {
@@ -228,13 +225,12 @@ async fn handle_server_session_tcp(
return;
}
};
- info!("{} how it went {:?}", sid ,stream);
// Signal success
mux.send(FrameType::ConnectOk, sid, Vec::new()).await;
// Data relay (shared implementation)
- info!("{} streaming", sid);
+ info!("[{}] streaming", sid);
relay_bidirectional_tcp(sid, stream, mux, session_rx).await;
info!("[{}] TCP Closed", sid);
} else {
@@ -250,7 +246,7 @@ async fn handle_server_session_tcp(
/// Handle a single proxied UDP session on the server side.
async fn handle_server_session_udp(
sid: u32,
- target_addr: TargetAddr,
+ _target_addr: TargetAddr,
mux: MuxHandle,
session_rx: mpsc::UnboundedReceiver<Frame>,
filter_config: FilterConfig
Served by rngit 1.4.2 - Generated in 0.05s